🎖️GitЯра🎖️
Node / meshtastic / Meshtastic-Android / files / core / takserver / src / commonMain / kotlin / org / meshtastic / core / takserver / TAKServerManager.kt
Displaying Raw • Download
core/takserver/src/commonMain/kotlin/org/meshtastic/core/takserver/TAKServerManager.kt bd2863243bab6eb213401d949839a2bc74dde7e2 (bd286324) Text, 9.04 KB
T8b949e/*
* Copyright (c) 2026 Meshtastic LLC
*
* This program is free software: you can redistribute it and/or modify
* it under the terms of the GNU General Public License as published by
* the Free Software Foundation, either version 3 of the License, or
* (at your option) any later version.
*
* This program is distributed in the hope that it will be useful,
* but WITHOUT ANY WARRANTY; without even the implied warranty of
* MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
* GNU General Public License for more details.
*
* You should have received a copy of the GNU General Public License
* along with this program. If not, see <https://www.gnu.org/licenses/>.
*/
Tff7b72package T7ee787org.meshtastic.core.takserver
Tff7b72import T7ee787co.touchlab.kermit.Logger
Tff7b72import T7ee787kotlinx.coroutines.CoroutineScope
Tff7b72import T7ee787kotlinx.coroutines.flow.MutableSharedFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.MutableStateFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.SharedFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.StateFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.asSharedFlow
Tff7b72import T7ee787kotlinx.coroutines.flow.asStateFlow
Tff7b72import T7ee787kotlinx.coroutines.launch
Tff7b72import T7ee787kotlinx.coroutines.sync.Mutex
Tff7b72import T7ee787kotlinx.coroutines.sync.withLock
Tff7b72import T7ee787kotlin.concurrent.Volatile
Tff7b72import T7ee787kotlin.time.Clock
Tff7b72import T7ee787kotlin.time.Duration.Companion.minutes
T8b949e/** A CoT message received from a connected TAK client, paired with the client's identity. */
Tff7b72data Tff7b72class T56d364InboundCoTMessageTb4b4b4(Tff7b72val Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4, Tff7b72val Te6edf3clientInfoTb4b4b4: Te6edf3TAKClientInfo? Tff7b72= Tff7b72nullTb4b4b4)
Tff7b72interface T56d364TAKServerManager Tb4b4b4{
Tff7b72val Te6edf3isSupportedTb4b4b4: Tffa657Boolean
Tff7b72val Te6edf3isRunningTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657BooleanTff7b72>
Tff7b72val Te6edf3isStartingTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657BooleanTff7b72>
Tff7b72val Te6edf3connectionCountTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657IntTff7b72>
Tff7b72val Te6edf3hasStartErrorTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657BooleanTff7b72>
Tff7b72val Te6edf3inboundMessagesTb4b4b4: Te6edf3SharedFlowTff7b72<Te6edf3InboundCoTMessageTff7b72>
T8b949e/**
* Emits once each time a TAK client connects.
*
* [TAKServer.onClientConnected] is a single callback slot already owned by the offline-queue drain, so consumers
* that need the connect event observe it here instead of replacing that callback.
*/
Tff7b72val Te6edf3clientConnectedTb4b4b4: Te6edf3SharedFlowTff7b72<Tffa657UnitTff7b72>
T8b949e/** Start the TAK server using [scope]. Port is fixed at [TAKServer] construction time. */
Tff7b72fun Td2a8ffstartTb4b4b4(Te6edf3scopeTb4b4b4: Te6edf3CoroutineScopeTb4b4b4)
Tff7b72fun Td2a8ffstopTb4b4b4(Tb4b4b4)
Tff7b72fun Td2a8ffbroadcastTb4b4b4(Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4)
T8b949e/** Broadcast raw XML verbatim to TAK clients, bypassing CoTMessage parsing. */
Tff7b72fun Td2a8ffbroadcastRawXmlTb4b4b4(Te6edf3xmlTb4b4b4: Tffa657StringTb4b4b4)
Tb4b4b4}
Tff7b72internal Tff7b72class T56d364TAKServerManagerImplTb4b4b4(Tff7b72private Tff7b72val Te6edf3takServerTb4b4b4: Te6edf3TAKServerTb4b4b4) Tb4b4b4: Te6edf3TAKServerManager Tb4b4b4{
Tff7b72private Tff7b72var Te6edf3scopeTb4b4b4: Te6edf3CoroutineScope? Tff7b72= Tff7b72null
Tf0883e@Volatile Tff7b72private Tff7b72var Te6edf3startGeneration Tff7b72= T79c0ff0L
Tff7b72override Tff7b72val Te6edf3isSupported Tff7b72= Te6edf3takServerTb4b4b4.Te6edf3isSupported
Tff7b72private Tff7b72val Te6edf3_isRunning Tff7b72= Te6edf3MutableStateFlowTb4b4b4(Tff7b72falseTb4b4b4)
Tff7b72override Tff7b72val Te6edf3isRunningTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657BooleanTff7b72> Tff7b72= Te6edf3_isRunningTb4b4b4.Te6edf3asStateFlowTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3_isStarting Tff7b72= Te6edf3MutableStateFlowTb4b4b4(Tff7b72falseTb4b4b4)
Tff7b72override Tff7b72val Te6edf3isStartingTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657BooleanTff7b72> Tff7b72= Te6edf3_isStartingTb4b4b4.Te6edf3asStateFlowTb4b4b4(Tb4b4b4)
T8b949e// Mirror TAKServer's event-driven connection count — no polling needed
Tff7b72override Tff7b72val Te6edf3connectionCountTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657IntTff7b72> Tff7b72= Te6edf3takServerTb4b4b4.Te6edf3connectionCount
Tff7b72private Tff7b72val Te6edf3_hasStartError Tff7b72= Te6edf3MutableStateFlowTb4b4b4(Tff7b72falseTb4b4b4)
Tff7b72override Tff7b72val Te6edf3hasStartErrorTb4b4b4: Te6edf3StateFlowTff7b72<Tffa657BooleanTff7b72> Tff7b72= Te6edf3_hasStartErrorTb4b4b4.Te6edf3asStateFlowTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3_inboundMessages Tff7b72= Te6edf3MutableSharedFlowTff7b72<Te6edf3InboundCoTMessageTff7b72>Tb4b4b4(Te6edf3extraBufferCapacity Tff7b72= T79c0ff6T79c0ff4Tb4b4b4)
Tff7b72override Tff7b72val Te6edf3inboundMessagesTb4b4b4: Te6edf3SharedFlowTff7b72<Te6edf3InboundCoTMessageTff7b72> Tff7b72= Te6edf3_inboundMessagesTb4b4b4.Te6edf3asSharedFlowTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3_clientConnected Tff7b72= Te6edf3MutableSharedFlowTff7b72<Tffa657UnitTff7b72>Tb4b4b4(Te6edf3extraBufferCapacity Tff7b72= T79c0ff8Tb4b4b4)
Tff7b72override Tff7b72val Te6edf3clientConnectedTb4b4b4: Te6edf3SharedFlowTff7b72<Tffa657UnitTff7b72> Tff7b72= Te6edf3_clientConnectedTb4b4b4.Te6edf3asSharedFlowTb4b4b4(Tb4b4b4)
T8b949e// Offline message queue — buffers mesh-originated CoT messages when no TAK
T8b949e// clients are connected, then drains them when a client reconnects. Entries
T8b949e// expire after OFFLINE_QUEUE_TTL to avoid delivering stale situational data.
Tff7b72private Tff7b72data Tff7b72class T56d364QueuedMessageTb4b4b4(Tff7b72val Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4, Tff7b72val Te6edf3enqueuedAtTb4b4b4: Te6edf3kotlinTb4b4b4.Te6edf3timeTb4b4b4.Te6edf3InstantTb4b4b4)
Tff7b72private Tff7b72val Te6edf3offlineQueue Tff7b72= Te6edf3ArrayDequeTff7b72<Te6edf3QueuedMessageTff7b72>Tb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3offlineQueueMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72private Tff7b72val Te6edf3lifecycleMutex Tff7b72= Te6edf3MutexTb4b4b4(Tb4b4b4)
Tff7b72companion Tff7b72object Tb4b4b4{
Tff7b72private Tff7b72val Te6edf3OFFLINE_QUEUE_TTL Tff7b72= T79c0ff5.Te6edf3minutes
Tff7b72private Tff7b72const Tff7b72val Te6edf3OFFLINE_QUEUE_MAX_SIZE Tff7b72= T79c0ff5T79c0ff0
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffstartTb4b4b4(Te6edf3scopeTb4b4b4: Te6edf3CoroutineScopeTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3isSupportedTb4b4b4) Tff7b72return
Tff7b72if Tb4b4b4(Te6edf3_isRunningTb4b4b4.Te6edf3value Tff7b72|Tff7b72| Te6edf3_isStartingTb4b4b4.Te6edf3valueTb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffTAKServerManager already runningTa5d6ff" Tb4b4b4}
Tff7b72return
Tb4b4b4}
Te6edf3_hasStartErrorTb4b4b4.Te6edf3value Tff7b72= Tff7b72false
Te6edf3_isStartingTb4b4b4.Te6edf3value Tff7b72= Tff7b72true
Tff7b72val Te6edf3generation Tff7b72= Tff7b72+Tff7b72+Te6edf3startGeneration
T8b949e// Assign scope AFTER the guard so a second concurrent start() can never
T8b949e// overwrite the active scope without actually restarting the server.
Tff7b72thisTb4b4b4.Te6edf3scope Tff7b72= Te6edf3scope
Te6edf3scopeTb4b4b4.Te6edf3launch Tb4b4b4{
Te6edf3lifecycleMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3generation Tff7b72!Tff7b72= Te6edf3startGenerationTb4b4b4) Tff7b72returnTf0883e@withLock
T8b949e// Wire up inbound message handler BEFORE starting so no messages are lost.
T8b949e// Use tryEmit (non-suspending) with extraBufferCapacity to avoid launching a
T8b949e// new coroutine per message, which would create unbounded coroutines under
T8b949e// high message rates and could reorder messages.
Te6edf3takServerTb4b4b4.Te6edf3onMessage Tff7b72= Tb4b4b4{ Te6edf3cotMessageTb4b4b4, Te6edf3clientInfo Tff7b72-Tff7b72>
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3_inboundMessagesTb4b4b4.Te6edf3tryEmitTb4b4b4(Te6edf3InboundCoTMessageTb4b4b4(Te6edf3cotMessageTb4b4b4, Te6edf3clientInfoTb4b4b4)Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3w Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK inbound message buffer full; dropping message from Tffd700${Te6edf3clientInfoTff7b72?.Te6edf3idTffd700}Ta5d6ff" Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Te6edf3takServerTb4b4b4.Te6edf3onClientConnected Tff7b72= Tb4b4b4{
Te6edf3drainOfflineQueueTb4b4b4(Tb4b4b4)
Te6edf3_clientConnectedTb4b4b4.Te6edf3tryEmitTb4b4b4(Tffa657UnitTb4b4b4)
Tb4b4b4}
Tff7b72val Te6edf3result Tff7b72= Te6edf3takServerTb4b4b4.Te6edf3startTb4b4b4(Te6edf3scopeTb4b4b4)
Tff7b72if Tb4b4b4(Te6edf3generation Tff7b72!Tff7b72= Te6edf3startGenerationTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3resultTb4b4b4.Te6edf3isSuccessTb4b4b4) Te6edf3takServerTb4b4b4.Te6edf3stopTb4b4b4(Tb4b4b4)
Tff7b72returnTf0883e@withLock
Tb4b4b4}
Te6edf3_isStartingTb4b4b4.Te6edf3value Tff7b72= Tff7b72false
Tff7b72if Tb4b4b4(Te6edf3resultTb4b4b4.Te6edf3isSuccessTb4b4b4) Tb4b4b4{
Te6edf3_isRunningTb4b4b4.Te6edf3value Tff7b72= Tff7b72true
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK Server startedTa5d6ff" Tb4b4b4}
Tb4b4b4} Tff7b72else Tb4b4b4{
Te6edf3_hasStartErrorTb4b4b4.Te6edf3value Tff7b72= Tff7b72true
Te6edf3LoggerTb4b4b4.Te6edf3eTb4b4b4(Te6edf3resultTb4b4b4.Te6edf3exceptionOrNullTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{ Ta5d6ff"Ta5d6ffFailed to start TAK ServerTa5d6ff" Tb4b4b4}
T8b949e// Clear both callbacks if start failed so we don't hold a reference unnecessarily
Te6edf3takServerTb4b4b4.Te6edf3onMessage Tff7b72= Tff7b72null
Te6edf3takServerTb4b4b4.Te6edf3onClientConnected Tff7b72= Tff7b72null
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffstopTb4b4b4(Tb4b4b4) Tb4b4b4{
T8b949e// Flip the running flag and null out the scope BEFORE stopping the server so
T8b949e// any broadcast()/drainOfflineQueue() that races stop() sees _isRunning=false
T8b949e// and exits early instead of launching coroutines on a scope that is about to
T8b949e// be discarded.
Te6edf3startGenerationTff7b72+Tff7b72+
Te6edf3_isRunningTb4b4b4.Te6edf3value Tff7b72= Tff7b72false
Te6edf3_isStartingTb4b4b4.Te6edf3value Tff7b72= Tff7b72false
Te6edf3_hasStartErrorTb4b4b4.Te6edf3value Tff7b72= Tff7b72false
Te6edf3scope Tff7b72= Tff7b72null
Te6edf3takServerTb4b4b4.Te6edf3onMessage Tff7b72= Tff7b72null
Te6edf3takServerTb4b4b4.Te6edf3onClientConnected Tff7b72= Tff7b72null
Te6edf3takServerTb4b4b4.Te6edf3stopTb4b4b4(Tb4b4b4)
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffTAK Server stoppedTa5d6ff" Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffbroadcastTb4b4b4(Te6edf3cotMessageTb4b4b4: Te6edf3CoTMessageTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3_isRunningTb4b4b4.Te6edf3valueTb4b4b4) Tff7b72return
Te6edf3scopeTff7b72?.Te6edf3launch Tb4b4b4{
Tff7b72if Tb4b4b4(Te6edf3takServerTb4b4b4.Te6edf3hasConnectionsTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3takServerTb4b4b4.Te6edf3broadcastTb4b4b4(Te6edf3cotMessageTb4b4b4)
Tb4b4b4} Tff7b72else Tb4b4b4{
T8b949e// No TAK clients connected — queue for delivery when one reconnects
Te6edf3offlineQueueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
T8b949e// Evict expired entries
Tff7b72val Te6edf3cutoff Tff7b72= Te6edf3ClockTb4b4b4.Te6edf3SystemTb4b4b4.Te6edf3nowTb4b4b4(Tb4b4b4) Tff7b72- Te6edf3OFFLINE_QUEUE_TTL
Tff7b72while Tb4b4b4(Te6edf3offlineQueueTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4) Tff7b72&Tff7b72& Te6edf3offlineQueueTb4b4b4.Te6edf3firstTb4b4b4(Tb4b4b4)Tb4b4b4.Te6edf3enqueuedAt Tff7b72< Te6edf3cutoffTb4b4b4) Tb4b4b4{
Te6edf3offlineQueueTb4b4b4.Te6edf3removeFirstTb4b4b4(Tb4b4b4)
Tb4b4b4}
T8b949e// Cap size to prevent unbounded growth
Tff7b72if Tb4b4b4(Te6edf3offlineQueueTb4b4b4.Te6edf3size Tff7b72>Tff7b72= Te6edf3OFFLINE_QUEUE_MAX_SIZETb4b4b4) Tb4b4b4{
Te6edf3offlineQueueTb4b4b4.Te6edf3removeFirstTb4b4b4(Tb4b4b4)
Tb4b4b4}
Te6edf3offlineQueueTb4b4b4.Te6edf3addLastTb4b4b4(Te6edf3QueuedMessageTb4b4b4(Te6edf3cotMessageTb4b4b4, Te6edf3ClockTb4b4b4.Te6edf3SystemTb4b4b4.Te6edf3nowTb4b4b4(Tb4b4b4)Tb4b4b4)Tb4b4b4)
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tff7b72override Tff7b72fun Td2a8ffbroadcastRawXmlTb4b4b4(Te6edf3xmlTb4b4b4: Tffa657StringTb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3_isRunningTb4b4b4.Te6edf3valueTb4b4b4) Tff7b72return
Te6edf3scopeTff7b72?.Te6edf3launch Tb4b4b4{ Te6edf3takServerTb4b4b4.Te6edf3broadcastRawXmlTb4b4b4(Te6edf3xmlTb4b4b4) Tb4b4b4}
Tb4b4b4}
T8b949e/**
* Drain any queued messages to the newly connected TAK client. Called by the server when a TAK client connects
* (Connected event).
*/
Tff7b72internal Tff7b72fun Td2a8ffdrainOfflineQueueTb4b4b4(Tb4b4b4) Tb4b4b4{
Tff7b72if Tb4b4b4(Tff7b72!Te6edf3_isRunningTb4b4b4.Te6edf3valueTb4b4b4) Tff7b72return
Te6edf3scopeTff7b72?.Te6edf3launch Tb4b4b4{
Tff7b72val Te6edf3messages Tff7b72=
Te6edf3offlineQueueMutexTb4b4b4.Te6edf3withLock Tb4b4b4{
Tff7b72val Te6edf3cutoff Tff7b72= Te6edf3ClockTb4b4b4.Te6edf3SystemTb4b4b4.Te6edf3nowTb4b4b4(Tb4b4b4) Tff7b72- Te6edf3OFFLINE_QUEUE_TTL
Tff7b72val Te6edf3valid Tff7b72= Te6edf3offlineQueueTb4b4b4.Te6edf3filter Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3enqueuedAt Tff7b72>Tff7b72= Te6edf3cutoff Tb4b4b4}Tb4b4b4.Te6edf3map Tb4b4b4{ Tffa657itTb4b4b4.Te6edf3cotMessage Tb4b4b4}
Te6edf3offlineQueueTb4b4b4.Te6edf3clearTb4b4b4(Tb4b4b4)
Te6edf3valid
Tb4b4b4}
Tff7b72if Tb4b4b4(Te6edf3messagesTb4b4b4.Te6edf3isNotEmptyTb4b4b4(Tb4b4b4)Tb4b4b4) Tb4b4b4{
Te6edf3LoggerTb4b4b4.Te6edf3i Tb4b4b4{ Ta5d6ff"Ta5d6ffDraining Tffd700${Te6edf3messagesTb4b4b4.Te6edf3sizeTffd700}Ta5d6ff queued message(s) to reconnected TAK clientTa5d6ff" Tb4b4b4}
Te6edf3messagesTb4b4b4.Te6edf3forEach Tb4b4b4{ Te6edf3takServerTb4b4b4.Te6edf3broadcastTb4b4b4(Tffa657itTb4b4b4) Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Tb4b4b4}
Served by rngit 1.5.0 - Generated in 0.17s